Skip to contentMethod: static {...}
1: /*
2: * *************************************************************************************************************************************************************
3: *
4: * TheseFoolishThings: Miscellaneous utilities
5: * http://tidalwave.it/projects/thesefoolishthings
6: *
7: * Copyright (C) 2009 - 2025 by Tidalwave s.a.s. (http://tidalwave.it)
8: *
9: * *************************************************************************************************************************************************************
10: *
11: * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License.
12: * You may obtain a copy of the License at
13: *
14: * http://www.apache.org/licenses/LICENSE-2.0
15: *
16: * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR
17: * CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License.
18: *
19: * *************************************************************************************************************************************************************
20: *
21: * git clone https://bitbucket.org/tidalwave/thesefoolishthings-src
22: * git clone https://github.com/tidalwave-it/thesefoolishthings-src
23: *
24: * *************************************************************************************************************************************************************
25: */
26: package it.tidalwave.messagebus.spi;
27:
28: import javax.annotation.Nonnull;
29: import it.tidalwave.messagebus.spi.MultiQueue.TopicAndMessage;
30: import lombok.Getter;
31: import lombok.Setter;
32: import lombok.ToString;
33: import lombok.extern.slf4j.Slf4j;
34:
35: /***************************************************************************************************************************************************************
36: *
37: * An implementation of {@link MessageDelivery} that dispatches messages in a round-robin fashion, topic by topic.
38: * Each delivery is performed in a separated thread.
39: *
40: * @author Fabrizio Giudici
41: * @since 2.2
42: *
43: **************************************************************************************************************************************************************/
44: @Slf4j @ToString(of = "workers")
45: public class RoundRobinAsyncMessageDelivery implements MessageDelivery
46: {
47: @Nonnull
48: private SimpleMessageBus messageBusSupport;
49:
50: @Getter @Setter
51: private int workers = 10;
52:
53: private final MultiQueue multiQueue = new MultiQueue();
54:
55: /***********************************************************************************************************************************************************
56: **********************************************************************************************************************************************************/
57: private final Runnable dispatcher = new Runnable()
58: {
59: @Override
60: public void run()
61: {
62: for (;;)
63: {
64: try
65: {
66: dispatchMessage(multiQueue.remove());
67: }
68: catch (InterruptedException e)
69: {
70: break;
71: }
72: }
73: }
74:
75: private <T> void dispatchMessage (@Nonnull final TopicAndMessage<T> tam)
76: {
77: messageBusSupport.dispatchMessage(tam.getTopic(), tam.getMessage());
78: }
79: };
80:
81: /***********************************************************************************************************************************************************
82: * {@inheritDoc}
83: **********************************************************************************************************************************************************/
84: @Override
85: public void initialize (@Nonnull final SimpleMessageBus messageBusSupport)
86: {
87: this.messageBusSupport = messageBusSupport;
88: final var executor = this.messageBusSupport.getExecutor();
89:
90: for (var i = 0; i < workers; i++)
91: {
92: executor.execute(dispatcher);
93: }
94: }
95:
96: /***********************************************************************************************************************************************************
97: * {@inheritDoc}
98: **********************************************************************************************************************************************************/
99: @Override
100: public <T> void deliverMessage (@Nonnull final Class<T> topic, @Nonnull final T message)
101: {
102: multiQueue.add(topic, message);
103: }
104: }